package ex.datastream.sources;

import org.apache.flink.api.common.serialization.SimpleStringSchema;
import org.apache.flink.connector.kafka.source.KafkaSource;
import org.apache.flink.connector.kafka.source.enumerator.initializer.OffsetsInitializer;

public class KafkaDataSource {
    public KafkaSource getKafkaSource(){
        // 加载数据源
        KafkaSource<String> source = KafkaSource.<String>builder()
                .setBootstrapServers("192.168.1.4:9092")
                .setTopics("test9")
                .setGroupId("my-group")
                .setStartingOffsets(OffsetsInitializer.latest())
                .setValueOnlyDeserializer(new SimpleStringSchema())
                .build();
        return source;
    }
}
